原文:Routing
状态:待校对
翻译:Bingjian-Zhu
校对:

CC-BY-SA

路由(Routing)

(使用Go客户端)

在前面的教程中,我们实现了一个简单的日志系统。可以把日志消息广播给多个接收者。

本篇教程中我们打算新增一个功能 —— 使得它能够只订阅消息的一个字集。例如,我们只需要把严重的错误日志信息写入日志文件(存储到磁盘),但同时仍然把所有的日志信息输出到控制台中

绑定(Bindings)

前面的例子,我们已经创建过绑定(bindings),代码如下:

  1. err = ch.QueueBind(
  2. q.Name, // queue name
  3. "", // routing key
  4. "logs", // exchange
  5. false,
  6. nil)

绑定(binding)是指交换机(exchange)和队列(queue)的关系。可以简单理解为:这个队列(queue)对这个交换机(exchange)的消息感兴趣。

绑定的时候可以带上一个额外的routing_key参数。为了避免与Channel.Publish的参数混淆,我们把它叫做绑定键binding key。以下是如何创建一个带绑定键的绑定。

  1. err = ch.QueueBind(
  2. q.Name, // queue name
  3. "black", // routing key
  4. "logs", // exchange
  5. false,
  6. nil)

绑定键的意义取决于交换机(exchange)的类型。我们之前使用过fanout 交换机会忽略这个值。

直连交换机(Direct exchange)

我们的日志系统广播所有的消息给所有的消费者(consumers)。我们打算扩展它,使其基于日志的严重程度进行消息过滤。例如我们也许只是希望将比较严重的错误(error)日志写入磁盘,以免在警告(warning)或者信息(info)日志上浪费磁盘空间。

我们使用的fanout 交换机没有足够的灵活性 —— 它能做的仅仅是广播。

我们将会使用direct 交换机来代替。路由的算法很简单 —— 交换机将会对binding keyrouting key进行精确匹配,从而确定消息该分发到哪个队列。

下图能够很好的描述这个场景:

路由(Routing) - 图2

在这个场景中,我们可以看到direct交换机 X和两个队列进行了绑定。第一个队列使用orange作为binding key,第二个队列有两个绑定,一个使用black作为binding key,另外一个使用green

这样以来,当消息发布到routing key为orange的交换机时,就会被路由到队列Q1。routing key为black或者green的消息就会路由到Q2。其他的所有消息都将会被丢弃。

多个绑定(Multiple bindings)

路由(Routing) - 图3

多个队列使用相同的binding key是合法的。这个例子中,我们可以添加一个X和Q1之间的绑定,使用black为binding key。这样一来,direct交换机就和fanout交换机的行为一样,会将消息广播到所有匹配的队列。带有routing key为black的消息会同时发送到Q1和Q2。

发送日志

我们将会发送消息到一个fanout,把日志等级作为routing key。这样接收日志就可以根据日志级别来选择它想要处理的日志。我们先看看发送日志。

我们需要创建一个交换机(exchange):

  1. err = ch.ExchangeDeclare(
  2. "logs_direct", // name
  3. "direct", // type
  4. true, // durable
  5. false, // auto-deleted
  6. false, // internal
  7. false, // no-wait
  8. nil, // arguments
  9. )

然后我们发送一则消息:

  1. err = ch.ExchangeDeclare(
  2. "logs_direct", // name
  3. "direct", // type
  4. true, // durable
  5. false, // auto-deleted
  6. false, // internal
  7. false, // no-wait
  8. nil, // arguments
  9. )
  10. failOnError(err, "Failed to declare an exchange")
  11. body := bodyFrom(os.Args)
  12. err = ch.Publish(
  13. "logs_direct", // exchange
  14. severityFrom(os.Args), // routing key
  15. false, // mandatory
  16. false, // immediate
  17. amqp.Publishing{
  18. ContentType: "text/plain",
  19. Body: []byte(body),
  20. })

我们假设日志等级的值是info、warning、error中的一个。

订阅

处理接收消息的方式和之前差不多,只有一个例外,我们将会为我们感兴趣的每个严重级别分别创建一个新的绑定。

  1. q, err := ch.QueueDeclare(
  2. "", // name
  3. false, // durable
  4. false, // delete when usused
  5. true, // exclusive
  6. false, // no-wait
  7. nil, // arguments
  8. )
  9. failOnError(err, "Failed to declare a queue")
  10. if len(os.Args) < 2 {
  11. log.Printf("Usage: %s [info] [warning] [error]", os.Args[0])
  12. os.Exit(0)
  13. }
  14. for _, s := range os.Args[1:] {
  15. log.Printf("Binding queue %s to exchange %s with routing key %s",
  16. q.Name, "logs_direct", s)
  17. err = ch.QueueBind(
  18. q.Name, // queue name
  19. s, // routing key
  20. "logs_direct", // exchange
  21. false,
  22. nil)
  23. failOnError(err, "Failed to bind a queue")
  24. }

代码整合

路由(Routing) - 图4

emit_log_direct.go的代码:

  1. package main
  2. import (
  3. "log"
  4. "os"
  5. "strings"
  6. "github.com/streadway/amqp"
  7. )
  8. func failOnError(err error, msg string) {
  9. if err != nil {
  10. log.Fatalf("%s: %s", msg, err)
  11. }
  12. }
  13. func main() {
  14. conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
  15. failOnError(err, "Failed to connect to RabbitMQ")
  16. defer conn.Close()
  17. ch, err := conn.Channel()
  18. failOnError(err, "Failed to open a channel")
  19. defer ch.Close()
  20. err = ch.ExchangeDeclare(
  21. "logs_direct", // name
  22. "direct", // type
  23. true, // durable
  24. false, // auto-deleted
  25. false, // internal
  26. false, // no-wait
  27. nil, // arguments
  28. )
  29. failOnError(err, "Failed to declare an exchange")
  30. body := bodyFrom(os.Args)
  31. err = ch.Publish(
  32. "logs_direct", // exchange
  33. severityFrom(os.Args), // routing key
  34. false, // mandatory
  35. false, // immediate
  36. amqp.Publishing{
  37. ContentType: "text/plain",
  38. Body: []byte(body),
  39. })
  40. failOnError(err, "Failed to publish a message")
  41. log.Printf(" [x] Sent %s", body)
  42. }
  43. func bodyFrom(args []string) string {
  44. var s string
  45. if (len(args) < 3) || os.Args[2] == "" {
  46. s = "hello"
  47. } else {
  48. s = strings.Join(args[2:], " ")
  49. }
  50. return s
  51. }
  52. func severityFrom(args []string) string {
  53. var s string
  54. if (len(args) < 2) || os.Args[1] == "" {
  55. s = "info"
  56. } else {
  57. s = os.Args[1]
  58. }
  59. return s
  60. }

receive_logs_direct.go的代码:

  1. package main
  2. import (
  3. "log"
  4. "os"
  5. "github.com/streadway/amqp"
  6. )
  7. func failOnError(err error, msg string) {
  8. if err != nil {
  9. log.Fatalf("%s: %s", msg, err)
  10. }
  11. }
  12. func main() {
  13. conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
  14. failOnError(err, "Failed to connect to RabbitMQ")
  15. defer conn.Close()
  16. ch, err := conn.Channel()
  17. failOnError(err, "Failed to open a channel")
  18. defer ch.Close()
  19. err = ch.ExchangeDeclare(
  20. "logs_direct", // name
  21. "direct", // type
  22. true, // durable
  23. false, // auto-deleted
  24. false, // internal
  25. false, // no-wait
  26. nil, // arguments
  27. )
  28. failOnError(err, "Failed to declare an exchange")
  29. q, err := ch.QueueDeclare(
  30. "", // name
  31. false, // durable
  32. false, // delete when usused
  33. true, // exclusive
  34. false, // no-wait
  35. nil, // arguments
  36. )
  37. failOnError(err, "Failed to declare a queue")
  38. if len(os.Args) < 2 {
  39. log.Printf("Usage: %s [info] [warning] [error]", os.Args[0])
  40. os.Exit(0)
  41. }
  42. for _, s := range os.Args[1:] {
  43. log.Printf("Binding queue %s to exchange %s with routing key %s",
  44. q.Name, "logs_direct", s)
  45. err = ch.QueueBind(
  46. q.Name, // queue name
  47. s, // routing key
  48. "logs_direct", // exchange
  49. false,
  50. nil)
  51. failOnError(err, "Failed to bind a queue")
  52. }
  53. msgs, err := ch.Consume(
  54. q.Name, // queue
  55. "", // consumer
  56. true, // auto ack
  57. false, // exclusive
  58. false, // no local
  59. false, // no wait
  60. nil, // args
  61. )
  62. failOnError(err, "Failed to register a consumer")
  63. forever := make(chan bool)
  64. go func() {
  65. for d := range msgs {
  66. log.Printf(" [x] %s", d.Body)
  67. }
  68. }()
  69. log.Printf(" [*] Waiting for logs. To exit press CTRL+C")
  70. <-forever
  71. }

如果你希望只是保存warning和error级别的日志到磁盘,只需要打开控制台并输入:

  1. $ go run receive_logs_direct.go warning error > logs_from_rabbit.log

如果你希望所有的日志信息都输出到屏幕中,打开一个新的终端,然后输入:

  1. go run receive_logs_direct.go info warning error
  2. # => [*] Waiting for logs. To exit press CTRL+C

如果要触发一个error级别的日志,只需要输入:

  1. go run emit_log_direct.go error "Run. Run. Or it will explode."
  2. # => [x] Sent 'error':'Run. Run. Or it will explode.'

这里是完整的代码:(emit_log_direct.goreceive_logs_direct.go)